NOTE

1.5 Redis Pipeline

What Redis pipeline is, why it reduces network I/O, usage examples, implementation notes, and comparison with M commands, Lua, and transactions.

Redis / CacheCreated Updated 1 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What Is Redis Pipeline?

  • Redis uses a client-server model and is a TCP server implementing a request/response protocol:
    • the client sends a query to the server and usually reads the server response from the socket in a blocking manner;
    • the server processes the command and sends the response back to the client.
  • One RTT (round-trip time) is roughly 250 ms. Even if the server can process 100K requests per second, only four requests per second can now be processed. In other words, all the time is spent on network I/O.
  • Pipeline means the client sends multiple commands to the server without waiting for replies, and reads the replies only in the final step.

2. Why Redis Pipeline Is Needed

Pipeline can reduce the number of network I/O operations.

Traditional pipeline
Time n network operations + n commands 1 network operation + n commands
Data volume n commands n commands

3. Redis Pipeline Usage

3.1. Bulk Insert a Large Amount of Data

  1. Prepare a text file
SET Key0 Value0
SET Key1 Value1
SET KeyN ValueN
  1. cat data.txt | redis-cli --pipe
cat data.txt | redis-cli --pipe
  1. Check on the server.

3.2. Lua with Pipeline

var seckilGoodsScriptHash atomic.String
var lock sync.Mutex

// BatchInitSeckillGoodsParticipant batch-initializes the number of participants for seckill goods
func (p *ActivitySeckillCardDao) BatchInitSeckillGoodsParticipant(ctx context.Context,
	expireTs int64, incrSeckills ...*model.IncrSeckillParticipant) error {
	scriptHash, err := p.getSeckilGoodsScriptHash(ctx)
	if err != nil {
		return err
	}

	pipeline, err := p.redisProxy.Pipeline(ctx)
	if err != nil {
		metrics.Counter("redis-[BatchInitSeckillGoodsParticipant]failed-to-get-Pipeline").Incr()
		log.ErrorContextf(ctx,
			"BatchInitSeckillGoodsParticipant: redis get pipeline failed. err=%v",
			err)
		return err
	}
	defer pipeline.Close()

	for _, incrSeckill := range incrSeckills {
		args := make([]interface{}, 0)
		activityKey := seckillParticipantsKey(incrSeckill.CardId,
			incrSeckill.ActivityId)
		goodsKey := seckillGoodsParticipantsKey(incrSeckill.ActivityGoodsId)
		args = append(args, scriptHash, 1, activityKey, goodsKey,
			incrSeckill.InitParticipants, expireTs)
		if err := pipeline.Send("EVALSHA", args...); err != nil {
			metrics.Counter("redis-[BatchInitSeckillGoodsParticipant]failed-to-Send-to-buffer").Incr()
			log.ErrorContextf(ctx,
				"BatchInitSeckillGoodsParticipant: redis send buffer failed. err=%v",
				err)
			return err
		}
	}
	if err := pipeline.Flush(); err != nil {
		metrics.Counter("redis-[BatchInitSeckillGoodsParticipant]Flush-failed").Incr()
		log.ErrorContextf(ctx,
			"BatchInitSeckillGoodsParticipant: redis flush to server failed. err=%v",
			err)
		return err
	}
	reply, err := pipeline.Receive()
	log.DebugContextf(ctx, "BatchInitSeckillGoodsParticipant: req=%v, rsp=%v",
		incrSeckills, reply)

	if err != nil {
		metrics.Counter("redis-[BatchInitSeckillGoodsParticipant]Receive-failed").Incr()
		log.ErrorContextf(ctx,
			"BatchInitSeckillGoodsParticipant: redis receive server reply failed. err=%v",
			err)
		return err
	}
	return nil
}

// Get the Lua script hash
func (p *ActivitySeckillCardDao) getSeckilGoodsScriptHash(ctx context.Context) (string,
	error) {
	hash := seckilGoodsScriptHash.Load()
	if hash != "" {
		return hash, nil
	}

	lock.Lock()
	defer lock.Unlock()
	hash = seckilGoodsScriptHash.Load()
	if hash != "" {
		return hash, nil
	}
	script := redis.NewScript(1, initSeckillGoodsParticipant)
	if err := script.Load(ctx, p.redisProxy); err != nil {
		metrics.Counter("redis-[BatchInitSeckillGoodsParticipant]failed-to-load-Script").Incr()
		log.ErrorContextf(ctx,
			"BatchInitSeckillGoodsParticipant: redis load script failed. err=%v",
			err)
		return "", err
	}

	seckilGoodsScriptHash.Store(script.Hash())
	return seckilGoodsScriptHash.Load(), nil
}

4. Redis Pipeline Implementation

When a client uses a pipeline to send commands, the server is forced to use memory to queue the replies. Therefore, if a large number of commands need to be sent through a pipeline, it is better to send them in batches.

5. Pipeline vs. M Commands vs. Lua Scripts vs. Transactions

M operation pipeline Lua script Transaction
Atomic? Yes No Yes Yes
Number of replies 1 n 1 1

6. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub